ποΈGitΠ―ΡΠ°ποΈ
Node / meshtastic / Meshtastic-Android / files / core / service / src / commonMain / kotlin / org / meshtastic / core / service / SharedRadioInterfaceService.kt
Displaying Raw β’ Download
core/service/src/commonMain/kotlin/org/meshtastic/core/service/SharedRadioInterfaceService.kt 01cd54907d62cd90913d3d071b6bc240cae365ef (01cd5490) Text, 58.61 KB
T8b949e/*
* Copyright (c) 2026 Meshtastic LLC
*
* This program is free software: you can redistribute it and/or modify
* it under the terms of the GNU General Public License as published by
* the Free Software Foundation, either version 3 of the License, or
* (at your option) any later version.
*
* This program is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
* GNU General Public License for more details.
*
* You should have received a copy of the GNU General Public License
* along with this program. If not, see <https://www.gnu.org/licenses/>.
*/
Tff7b72package T7ee787org.meshtastic.core.service
Tff7b72import T7ee787androidx.lifecycle.Lifecycle
Tff7b72import T7ee787androidx.lifecycle.coroutineScope
Tff7b72import T7ee787co.touchlab.kermit.Logger
Tff7b72import T7ee787kotlinx.atomicfu.atomic
Tff7b72import T7ee787kotlinx.atomicfu.locks.SynchronizedObject
Tff7b72import T7ee787kotlinx.atomicfu.locks.synchronized
Tff7b72import T7ee787kotlinx.coroutines.CompletableDeferred
Tff7b72import T7ee787kotlinx.coroutines.CoroutineScope
Tff7b72import T7ee787kotlinx.coroutines.Job
Tff7b72import T7ee787kotlinx.coroutines.NonCancellable
Tff7b72import T7ee787kotlinx.coroutines.SupervisorJob
Tff7b72import T7ee787kotlinx.coroutines.cancel
Tff7b72import T7ee787kotlinx.coroutines.channels.BufferOverflow
Tff7b72import T7ee787kotlinx.coroutines.channels.Channel
Tff7b72import T7ee787kotlinx.coroutines.delay
Tff7b72import T7ee787kotlinx.coroutines.flow.Flow
Tff7b72import T7ee787kotlinx.coroutines.flow.MutableSharedFlow
Tff7b72import T7ee787kotlinx.coroutines.flow.MutableStateFlow
Tff7b72import T7ee787kotlinx.coroutines.flow.StateFlow
Tff7b72import T7ee787kotlinx.coroutines.flow.asFlow
Tff7b72import T7ee787kotlinx.coroutines.flow.asStateFlow
Tff7b72import T7ee787kotlinx.coroutines.flow.catch
Tff7b72import T7ee787kotlinx.coroutines.flow.combine
Tff7b72import T7ee787kotlinx.coroutines.flow.distinctUntilChanged
Tff7b72import T7ee787kotlinx.coroutines.flow.drop
Tff7b72import T7ee787kotlinx.coroutines.flow.filter
Tff7b72import T7ee787kotlinx.coroutines.flow.launchIn
Tff7b72import T7ee787kotlinx.coroutines.flow.onEach
Tff7b72import T7ee787kotlinx.coroutines.flow.receiveAsFlow
Tff7b72import T7ee787kotlinx.coroutines.flow.runningFold
Tff7b72import T7ee787kotlinx.coroutines.launch
Tff7b72import T7ee787kotlinx.coroutines.sync.Mutex
Tff7b72import T7ee787kotlinx.coroutines.sync.withLock
Tff7b72import T7ee787kotlinx.coroutines.withContext
Tff7b72import T7ee787okio.ByteString.Companion.toByteString
Tff7b72import T7ee787org.koin.core.annotation.Named
Tff7b72import T7ee787org.koin.core.annotation.Single
Tff7b72import T7ee787org.meshtastic.core.ble.BluetoothRepository
Tff7b72import T7ee787org.meshtastic.core.common.util.handledLaunch
Tff7b72import T7ee787org.meshtastic.core.common.util.ignoreExceptionSuspend
Tff7b72import T7ee787org.meshtastic.core.common.util.nowMillis
Tff7b72import T7ee787org.meshtastic.core.di.CoroutineDispatchers
Tff7b72import T7ee787org.meshtastic.core.model.ConnectionState
Tff7b72import T7ee787org.meshtastic.core.model.DeviceType
Tff7b72import T7ee787org.meshtastic.core.model.InterfaceId
Tff7b72import T7ee787org.meshtastic.core.model.MeshActivity
Tff7b72import T7ee787org.meshtastic.core.model.util.anonymize
Tff7b72import T7ee787org.meshtastic.core.network.repository.NetworkRepository
Tff7b72import T7ee787org.meshtastic.core.network.repository.SerialDevicePresence
Tff7b72import T7ee787org.meshtastic.core.repository.PlatformAnalytics
Tff7b72import T7ee787org.meshtastic.core.repository.RadioInterfaceService
Tff7b72import T7ee787org.meshtastic.core.repository.RadioPrefs
Tff7b72import T7ee787org.meshtastic.core.repository.RadioSessionContext
Tff7b72import T7ee787org.meshtastic.core.repository.RadioSessionLease
Tff7b72import T7ee787org.meshtastic.core.repository.RadioTransport
Tff7b72import T7ee787org.meshtastic.core.repository.RadioTransportFactory
Tff7b72import T7ee787org.meshtastic.core.repository.ReceivedRadioFrame
Tff7b72import T7ee787org.meshtastic.core.repository.TransportDisconnectReason
Tff7b72import T7ee787org.meshtastic.proto.ToRadio
Tff7b72import T7ee787kotlin.concurrent.Volatile
Tff7b72private Tff7b72const Tff7b72val Te6edf3USB_PERMISSION_DENIED_ERROR Tff7b72= Ta5d6ff"Ta5d6ffUSB permission denied. Reconnect the device to try again.Ta5d6ff"
T8b949e/**
* Immutable per-start transport identity. [generation] is bumped on every transport start (including same-address
* reconnect) so the [SharedRadioInterfaceService] and its consumers can discard state retained from a previous
* transport instance even when the selected address is unchanged.
*/
Tff7b72data Tff7b72class T56d364RadioTransportSessionTb4b4b4(Tff7b72val Te6edf3generationTb4b4b4: Tffa657LongTb4b4b4, Tff7b72val Te6edf3addressTb4b4b4: Tffa657StringTb4b4b4) Tb4b4b4{
T8b949e/** Immutable context carried with every frame admitted by this session. */
Tff7b72val Te6edf3context Tff7b72= Te6edf3RadioSessionContextTb4b4b4(Te6edf3generation Tff7b72= Te6edf3generationTb4b4b4, Te6edf3address Tff7b72= Te6edf3addressTb4b4b4)
T8b949e/** Prevent accidental disclosure of the raw transport address in diagnostic interpolation. */
Tff7b72override Tff7b72fun Td2a8fftoStringTb4b4b4(Tb4b4b4)Tb4b4b4: Tffa657String Tff7b72= Ta5d6ff"Ta5d6ffRadioTransportSession(generation=Tffd700$Te6edf3generationTa5d6ff, address=...)Ta5d6ff"
Tb4b4b4}
Tff7b72private Tff7b72data Tff7b72class T56d364SelectedSerialPresenceTb4b4b4(Tff7b72val Te6edf3keyTb4b4b4: Tffa657String?Tb4b4b4, Tff7b72val Te6edf3presentTb4b4b4: Tffa657BooleanTb4b4b4)
Tff7b72private Tff7b72data Tff7b72class T56d364UsbRecoverySnapshotTb4b4b4(Tff7b72val Te6edf3presenceTb4b4b4: Te6edf3SelectedSerialPresenceTb4b4b4, Tff7b72val Te6edf3stateTb4b4b4: Te6edf3ConnectionStateTb4b4b4)
Tff7b72private Tff7b72data Tff7b72class T56d364UsbRecoveryTriggerStateTb4b4b4(
Tff7b72val Te6edf3keyTb4b4b4: Tffa657String? Tff7b72= Tff7b72nullTb4b4b4,
Tff7b72val Te6edf3presentTb4b4b4: Tffa657Boolean Tff7b72= Tff7b72falseTb4b4b4,
Tff7b72val Te6edf3armedByAbsenceTb4b4b4: Tffa657Boolean Tff7b72= Tff7b72falseTb4b4b4,
Tff7b72val Te6edf3triggerRecoveryTb4b4b4: Tffa657Boolean Tff7b72= Tff7b72falseTb4b4b4,
Tb4b4b4)
Tff7b72private Tff7b72fun Td2a8ffselectedSerialPresenceTb4b4b4(Te6edf3addressTb4b4b4: Tffa657String?Tb4b4b4, Te6edf3keysTb4b4b4: Te6edf3SetTff7b72<Tffa657StringTff7b72>Tb4b4b4)Tb4b4b4: Te6edf3SelectedSerialPresence Tb4b4b4{
Tff7b72val Te6edf3key Tff7b72= Te6edf3addressTff7b72?.Te6edf3takeIf Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3firstOrNullTb4b4b4(Tb4b4b4) Tff7b72=Tff7b72= Te6edf3InterfaceIdTb4b4b4.Te6edf3SERIALTb4b4b4.Te6edf3id Tb4b4b4}Tff7b72?.Te6edf3dropTb4b4b4(T79c0ff1Tb4b4b4)Tff7b72?.Te6edf3takeIf Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3isNotEmptyTb4b4b4(Tb4b4b4) Tb4b4b4}
Tff7b72return Te6edf3SelectedSerialPresenceTb4b4b4(Te6edf3key Tff7b72= Te6edf3keyTb4b4b4, Te6edf3present Tff7b72= Te6edf3key Tff7b72!Tff7b72= Tff7b72null Tff7b72&Tff7b72& Te6edf3key Tff7b72in Te6edf3keysTb4b4b4)
Tb4b4b4}
Tff7b72private Tff7b72fun Te6edf3UsbRecoveryTriggerStateTb4b4b4.Td2a8ffnextTb4b4b4(Te6edf3snapshotTb4b4b4: Te6edf3UsbRecoverySnapshotTb4b4b4)Tb4b4b4: Te6edf3UsbRecoveryTriggerState Tb4b4b4{
Tff7b72val Te6edf3key Tff7b72= Te6edf3snapshotTb4b4b4.Te6edf3presenceTb4b4b4.Te6edf3key
Tff7b72val Te6edf3present Tff7b72= Te6edf3snapshotTb4b4b4.Te6edf3presenceTb4b4b4.Te6edf3present
Tff7b72return Tff7b72when Tb4b4b4{
Te6edf3key Tff7b72=Tff7b72= Tff7b72null Tff7b72-Tff7b72> Te6edf3UsbRecoveryTriggerStateTb4b4b4(Tb4b4b4)
Te6edf3snapshotTb4b4b4.Te6edf3state Tff7b72=Tff7b72= Te6edf3ConnectionStateTb4b4b4.Te6edf3Disconnected Tff7b72-Tff7b72> Te6edf3UsbRecoveryTriggerStateTb4b4b4(Te6edf3key Tff7b72= Te6edf3keyTb4b4b4, Te6edf3present Tff7b72= Te6edf3presentTb4b4b4)
Te6edf3key Tff7b72!Tff7b72= Tff7b72thisTb4b4b4.Te6edf3key Tff7b72-Tff7b72>
Te6edf3UsbRecoveryTriggerStateTb4b4b4(
Te6edf3key Tff7b72= Te6edf3keyTb4b4b4,
Te6edf3present Tff7b72= Te6edf3presentTb4b4b4,
Te6edf3armedByAbsence Tff7b72= Tff7b72!Te6edf3present Tff7b72&Tff7b72& Te6edf3snapshotTb4b4b4.Te6edf3state Tff7b72=Tff7b72= Te6edf3ConnectionStateTb4b4b4.Te6edf3DeviceSleepTb4b4b4,
Tb4b4b4)
Tff7b72!Te6edf3present Tff7b72-Tff7b72>
Te6edf3UsbRecoveryTriggerStateTb4b4b4(
Te6edf3key Tff7b72= Te6edf3keyTb4b4b4,
Te6edf3armedByAbsence Tff7b72= Te6edf3armedByAbsence Tff7b72|Tff7b72| Tff7b72thisTb4b4b4.Te6edf3present Tff7b72|Tff7b72| Te6edf3snapshotTb4b4b4.Te6edf3state Tff7b72=Tff7b72= Te6edf3ConnectionStateTb4b4b4.Te6edf3DeviceSleepTb4b4b4,
Tb4b4b4)
Tff7b72else Tff7b72-Tff7b72> Tb4b4b4{
Tff7b72val Te6edf3trigger Tff7b72= Te6edf3armedByAbsence Tff7b72&Tff7b72& Te6edf3snapshotTb4b4b4.Te6edf3state Tff7b72=Tff7b72= Te6edf3ConnectionStateTb4b4b4.Te6edf3DeviceSleep
Te6edf3UsbRecoveryTriggerStateTb4b4b4(
Te6edf3key Tff7b72= Te6edf3keyTb4b4b4,
Te6edf3present Tff7b72= Tff7b72trueTb4b4b4,
Te6edf3armedByAbsence Tff7b72= Te6edf3armedByAbsence Tff7b72&Tff7b72& Tff7b72!Te6edf3triggerTb4b4b4,
Te6edf3triggerRecovery Tff7b72= Te6edf3triggerTb4b4b4,
Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72private Tff7b72fun Te6edf3TransportDisconnectReasonTb4b4b4.Td2a8fftoConnectionErrorMessageTb4b4b4(Tb4b4b4)Tb4b4b4: Tffa657String Tff7b72= Tff7b72when Tb4b4b4(Tff7b72thisTb4b4b4) Tb4b4b4{
Te6edf3TransportDisconnectReasonTb4b4b4.Te6edf3UsbPermissionDenied Tff7b72-Tff7b72> Te6edf3USB_PERMISSION_DENIED_ERROR
Tb4b4b4}
T8b949e/**
* Shared multiplatform connection orchestrator for Meshtastic radios.
*
* Manages the connection lifecycle (connect, active, disconnect, reconnect loop), device address state flows, and
* hardware state observability (BLE/Network toggles). Delegates the actual raw byte transport mapping to a
* platform-specific [RadioTransportFactory].
*/
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffLongParameterListTa5d6ff"Tb4b4b4, Ta5d6ff"Ta5d6ffTooManyFunctionsTa5d6ff"Tb4b4b4)
Tf0883e@Single
Tff7b72class T56d364SharedRadioInterfaceServiceTb4b4b4(
Tff7b72private Tff7b72val Te6edf3dispatchersTb4b4b4: Te6edf3CoroutineDispatchersTb4b4b4,
Tff7b72private Tff7b72val Te6edf3bluetoothRepositoryTb4b4b4: Te6edf3BluetoothRepositoryTb4b4b4,
Tff7b72private Tff7b72val Te6edf3networkRepositoryTb4b4b4: Te6edf3NetworkRepositoryTb4b4b4,
Tff7b72private Tff7b72val Te6edf3serialDevicePresenceTb4b4b4: Te6edf3SerialDevicePresenceTb4b4b4,
Tf0883e@NamedTb4b4b4(Ta5d6ff"Ta5d6ffProcessLifecycleTa5d6ff"Tb4b4b4) Tff7b72private Tff7b72val Te6edf3processLifecycleTb4b4b4: Te6edf3LifecycleTb4b4b4,
Tff7b72private Tff7b72val Te6edf3radioPrefsTb4b4b4: Te6edf3RadioPrefsTb4b4b4,
Tff7b72private Tff7b72val Te6edf3transportFactoryTb4b4b4: Te6edf3RadioTransportFactoryTb4b4b4,
Tff7b72private Tff7b72val Te6edf3analyticsTb4b4b4: Te6edf3PlatformAnalyticsTb4b4b4,
Tb4b4b4) Tb4b4b4: Te6edf3RadioInterfaceService Tb4b4b4{
Tff7b72override Tff7b72val Te6edf3supportedDeviceTypesTb4b4b4: Te6edf3ListTff7b72<Te6edf3DeviceTypeTff7b72>
Tff7b72getTb4b4b4(Tb4b4b4) Tff7b72= Te6edf3transportFactoryTb4b4b4.Te6edf3supportedDeviceTypes
T8b949e/**
* Transport-level connection state reflecting the raw hardware link status.
*
* Updated directly by [onConnect] and [onDisconnect] when the physical transport (BLE, TCP, Serial) connects or
* disconnects. This is consumed exclusively by
* [MeshConnectionManager][org.meshtastic.core.repository.MeshConnectionManager], which reconciles it into the
* canonical app-level
* [ServiceRepository.connectionState][org.meshtastic.core.repository.ServiceRepository.connectionState].
*/
Tff7b72private Tff7b72val Te6edf3_connectionState Tff7b72= Te6edf3MutableStateFlowTff7b72<Te6edf3ConnectionStateTff7b72>Tb4b4b4(Te6edf3ConnectionStateTb4b4b4.Te6edf3DisconnectedTb4b4b4)
Tff7b72override Tff7b72val Te6edf3connectionStateTb4b4b4: Te6edf3StateFlowTff7b72<Te6edf3ConnectionStateTff7b72> Tff7b72= Te6edf3_connectionStateTb4b4b4.Te6edf3asStateFlowTb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3_currentDeviceAddressFlow Tff7b72= Te6edf3MutableStateFlowTff7b72<Tffa657String?Tff7b72>Tb4b4b4(Te6edf3radioPrefsTb4b4b4.Te6edf3devAddrTb4b4b4.Te6edf3valueTb4b4b4)
Tff7b72override Tff7b72val Te6edf3currentDeviceAddressFlowTb4b4b4: Te6edf3StateFlowTff7b72<Tffa657String?Tff7b72> Tff7b72= Te6edf3_currentDeviceAddressFlowTb4b4b4.Te6edf3asStateFlowTb4b4b4(Tb4b4b4)
T8b949e// Monotonically increasing generation bumped on every transport start (including same-address reconnect). Exposed
T8b949e// through [sessionGeneration] so the controller layer can clear connection-session identity at each session
T8b949e// boundary instead of only when the selected address changes.
Tff7b72private Tff7b72val Te6edf3sessionGenerationCounter Tff7b72= Te6edf3atomicTb4b4b4(T79c0ff0LTb4b4b4)
Tff7b72private Tff7b72val Te6edf3_sessionGeneration Tff7b72= Te6edf3MutableStateFlowTb4b4b4(T79c0ff0LTb4b4b4)
Tff7b72override Tff7b72val Te6edf3sessionGenerationTb4b4b4: Te6edf3StateFlowTff7b72<Tffa657LongTff7b72> Tff7b72= Te6edf3_sessionGenerationTb4b4b4.Te6edf3asStateFlowTb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3_activeSession Tff7b72= Te6edf3MutableStateFlowTff7b72<Te6edf3RadioSessionContext?Tff7b72>Tb4b4b4(Tff7b72nullTb4b4b4)
Tff7b72override Tff7b72val Te6edf3activeSessionTb4b4b4: Te6edf3StateFlowTff7b72<Te6edf3RadioSessionContext?Tff7b72> Tff7b72= Te6edf3_activeSessionTb4b4b4.Te6edf3asStateFlowTb4b4b4(Tb4b4b4)
T8b949e/**
* The transport session that owns lifecycle completion. [stopTransportLocked] closes its admission first, waits for
* every already-admitted operation, then clears this token before closing the underlying transport. Close-time and
* post-close callbacks are rejected as soon as admission closes.
*/
Tf0883e@Volatile Tff7b72private Tff7b72var Te6edf3activeTransportSessionTb4b4b4: Te6edf3RadioTransportSession? Tff7b72= Tff7b72null
T8b949e/** Serializes session publication/revocation with callback and suspend-operation admission. */
Tff7b72private Tff7b72val Te6edf3sessionCallbackLock Tff7b72= Te6edf3SynchronizedObjectTb4b4b4(Tb4b4b4)
T8b949e/** Guarded by [sessionCallbackLock]. New work is rejected immediately when teardown closes this gate. */
Tff7b72private Tff7b72var Te6edf3sessionAdmissionOpen Tff7b72= Tff7b72false
T8b949e/** Number of suspend operations admitted for [activeTransportSession], guarded by [sessionCallbackLock]. */
Tff7b72private Tff7b72var Te6edf3admittedSessionOperations Tff7b72= T79c0ff0
T8b949e/** Completed by the last admitted operation after teardown closes admission. Guarded by [sessionCallbackLock]. */
Tff7b72private Tff7b72var Te6edf3sessionDrainWaiterTb4b4b4: Te6edf3CompletableDeferredTff7b72<Tffa657UnitTff7b72>Tff7b72? Tff7b72= Tff7b72null
T8b949e/** Preserves FIFO ordering for handshake work without blocking independently leased packet side effects. */
Tff7b72private Tff7b72val Te6edf3sessionOperationMutex Tff7b72= Te6edf3MutexTb4b4b4(Tb4b4b4)
Tff7b72override Tff7b72fun Td2a8ffisSessionActiveTb4b4b4(Te6edf3sessionTb4b4b4: Te6edf3RadioSessionContextTb4b4b4)Tb4b4b4: Tffa657Boolean Tff7b72=
Te6edf3synchronizedTb4b4b4(Te6edf3sessionCallbackLockTb4b4b4) Tb4b4b4{ Te6edf3sessionAdmissionOpen Tff7b72&Tff7b72& Te6edf3activeTransportSessionTff7b72?.Te6edf3context Tff7b72=Tff7b72= Te6edf3session Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffrunIfSessionActiveTb4b4b4(Te6edf3sessionTb4b4b4: Te6edf3RadioSessionContextTb4b4b4, Te6edf3blockTb4b4b4: Tb4b4b4(Tb4b4b4) Tff7b72-Tff7b72> Tffa657UnitTb4b4b4)Tb4b4b4: Tffa657Boolean Tff7b72=
Te6edf3synchronizedTb4b4b4(Te6edf3sessionCallbackLockTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3sessionAdmissionOpen Tff7b72|Tff7b72| Te6edf3activeTransportSessionTff7b72?.Te6edf3context Tff7b72!Tff7b72= Te6edf3sessionTb4b4b4) Tff7b72returnTf0883e@synchronized Tff7b72false
Te6edf3blockTb4b4b4(Tb4b4b4)
Tff7b72true
Tb4b4b4}
Tff7b72override Tff7b72suspend Tff7b72fun Td2a8ffrunWithSessionLeaseTb4b4b4(
Te6edf3sessionTb4b4b4: Te6edf3RadioSessionContextTb4b4b4,
Te6edf3blockTb4b4b4: Te6edf3suspend Tb4b4b4(Te6edf3RadioSessionLeaseTb4b4b4) Tff7b72-Tff7b72> Tffa657UnitTb4b4b4,
Tb4b4b4)Tb4b4b4: Tffa657Boolean Tb4b4b4{
Tff7b72val Te6edf3admittedSession Tff7b72=
Te6edf3synchronizedTb4b4b4(Te6edf3sessionCallbackLockTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3active Tff7b72= Te6edf3activeTransportSession
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3sessionAdmissionOpen Tff7b72|Tff7b72| Te6edf3activeTff7b72?.Te6edf3context Tff7b72!Tff7b72= Te6edf3sessionTb4b4b4) Tb4b4b4{
Tff7b72null
Tb4b4b4} Tff7b72else Tb4b4b4{
Te6edf3admittedSessionOperationsTff7b72+Tff7b72+
Te6edf3active
Tb4b4b4}
Tb4b4b4} Tff7b72?: Tff7b72return Tff7b72false
Tff7b72val Te6edf3lease Tff7b72=
Tff7b72object Tb4b4b4: T56d364RadioSessionLease Tb4b4b4{
Tff7b72override Tff7b72val Te6edf3sessionTb4b4b4: Te6edf3RadioSessionContext Tff7b72= Te6edf3session
Tff7b72override Tff7b72fun Td2a8ffisCurrentTb4b4b4(Tb4b4b4)Tb4b4b4: Tffa657Boolean Tff7b72=
Te6edf3synchronizedTb4b4b4(Te6edf3sessionCallbackLockTb4b4b4) Tb4b4b4{ Te6edf3activeTransportSession Tff7b72=Tff7b72=Tff7b72= Te6edf3admittedSession Tb4b4b4}
Tb4b4b4}
Tff7b72try Tb4b4b4{
Te6edf3blockTb4b4b4(Te6edf3leaseTb4b4b4)
Tff7b72return Tff7b72true
Tb4b4b4} Tff7b72finally Tb4b4b4{
Tff7b72val Te6edf3drainWaiter Tff7b72=
Te6edf3synchronizedTb4b4b4(Te6edf3sessionCallbackLockTb4b4b4) Tb4b4b4{
Te6edf3checkTb4b4b4(Te6edf3activeTransportSession Tff7b72=Tff7b72=Tff7b72= Te6edf3admittedSessionTb4b4b4) Tb4b4b4{
Ta5d6ff"Ta5d6ffSession changed before an admitted operation released its leaseTa5d6ff"
Tb4b4b4}
Te6edf3checkTb4b4b4(Te6edf3admittedSessionOperations Tff7b72> T79c0ff0Tb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffSession operation count underflowTa5d6ff" Tb4b4b4}
Te6edf3admittedSessionOperationsTff7b72-Tff7b72-
Tff7b72if Tb4b4b4(Te6edf3admittedSessionOperations Tff7b72=Tff7b72= T79c0ff0Tb4b4b4) Tb4b4b4{
Te6edf3sessionDrainWaiterTb4b4b4.Te6edf3also Tb4b4b4{ Te6edf3sessionDrainWaiter Tff7b72= Tff7b72null Tb4b4b4}
Tb4b4b4} Tff7b72else Tb4b4b4{
Tff7b72null
Tb4b4b4}
Tb4b4b4}
Te6edf3drainWaiterTff7b72?.Te6edf3completeTb4b4b4(Tffa657UnitTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72suspend Tff7b72fun Td2a8ffrunWhileSessionActiveTb4b4b4(Te6edf3sessionTb4b4b4: Te6edf3RadioSessionContextTb4b4b4, Te6edf3blockTb4b4b4: Te6edf3suspend Tb4b4b4(Tb4b4b4) Tff7b72-Tff7b72> Tffa657UnitTb4b4b4)Tb4b4b4: Tffa657Boolean Tff7b72=
Te6edf3sessionOperationMutexTb4b4b4.Te6edf3withLock Tb4b4b4{ Te6edf3runWithSessionLeaseTb4b4b4(Te6edf3sessionTb4b4b4) Tb4b4b4{ Te6edf3blockTb4b4b4(Tb4b4b4) Tb4b4b4} Tb4b4b4}
T8b949e/** Runs a callback only while [session] still owns admission, atomically with session teardown. */
Tff7b72private Tff7b72inline Tff7b72fun Td2a8ffrunIfTransportSessionActiveTb4b4b4(Te6edf3sessionTb4b4b4: Te6edf3RadioTransportSessionTb4b4b4, Te6edf3blockTb4b4b4: Tb4b4b4(Tb4b4b4) Tff7b72-Tff7b72> Tffa657UnitTb4b4b4)Tb4b4b4: Tffa657Boolean Tff7b72=
Te6edf3synchronizedTb4b4b4(Te6edf3sessionCallbackLockTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3sessionAdmissionOpen Tff7b72|Tff7b72| Te6edf3activeTransportSession Tff7b72!Tff7b72=Tff7b72= Te6edf3sessionTb4b4b4) Tff7b72returnTf0883e@synchronized Tff7b72false
Te6edf3blockTb4b4b4(Tb4b4b4)
Tff7b72true
Tb4b4b4}
T8b949e/**
* Closes admission for [session], waits for its existing suspend operations, then publishes lifecycle completion.
* The closed gate rejects queued callbacks and new operations immediately; the retained session token prevents a
* replacement generation from overlapping an admitted database commit.
*/
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffrevokeTransportSessionTb4b4b4(Te6edf3sessionTb4b4b4: Te6edf3RadioTransportSession?Tb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3session Tff7b72=Tff7b72= Tff7b72nullTb4b4b4) Tff7b72return
Te6edf3withContextTb4b4b4(Te6edf3NonCancellableTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3drainWaiter Tff7b72=
Te6edf3synchronizedTb4b4b4(Te6edf3sessionCallbackLockTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3activeTransportSession Tff7b72!Tff7b72=Tff7b72= Te6edf3sessionTb4b4b4) Tff7b72returnTf0883e@synchronized Tff7b72null
Te6edf3sessionAdmissionOpen Tff7b72= Tff7b72false
Tff7b72if Tb4b4b4(Te6edf3admittedSessionOperations Tff7b72=Tff7b72= T79c0ff0Tb4b4b4) Tb4b4b4{
Tff7b72null
Tb4b4b4} Tff7b72else Tb4b4b4{
Te6edf3sessionDrainWaiter Tff7b72?: Te6edf3CompletableDeferredTff7b72<Tffa657UnitTff7b72>Tb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3also Tb4b4b4{ Te6edf3sessionDrainWaiter Tff7b72= Tffa657it Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Te6edf3drainWaiterTff7b72?.Te6edf3awaitTb4b4b4(Tb4b4b4)
Te6edf3synchronizedTb4b4b4(Te6edf3sessionCallbackLockTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3activeTransportSession Tff7b72=Tff7b72=Tff7b72= Te6edf3sessionTb4b4b4) Tb4b4b4{
Te6edf3checkTb4b4b4(Te6edf3admittedSessionOperations Tff7b72=Tff7b72= T79c0ff0Tb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffSession revoked before admitted operations drainedTa5d6ff" Tb4b4b4}
Te6edf3activeTransportSession Tff7b72= Tff7b72null
Te6edf3_activeSessionTb4b4b4.Te6edf3value Tff7b72= Tff7b72null
Te6edf3sessionDrainWaiter Tff7b72= Tff7b72null
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
T8b949e// A Channel preserves strict FIFO delivery of incoming radio bytes, which the firmware
T8b949e// handshake depends on (initial config packet ordering). A SharedFlow with `launch { emit() }`
T8b949e// per packet reorders under concurrent dispatch and breaks config load. trySend never
T8b949e// suspends, so handleFromRadio can remain a non-suspend synchronous callback.
T8b949e//
T8b949e// Bounded rather than UNLIMITED: inbound frames arrive faster than they can be processed under
T8b949e// sustained traffic, and an unbounded queue grows with no drop policy at all. At the cap
T8b949e// trySend fails and the newest frame is dropped, which keeps the already-queued (earlier)
T8b949e// frames and so preserves the ordering the handshake relies on.
Tff7b72private Tff7b72val Te6edf3_receivedData Tff7b72= Te6edf3ChannelTff7b72<Te6edf3ReceivedRadioFrameTff7b72>Tb4b4b4(Te6edf3RECEIVE_QUEUE_CAPACITYTb4b4b4)
Tff7b72override Tff7b72val Te6edf3receivedDataTb4b4b4: Te6edf3FlowTff7b72<Te6edf3ReceivedRadioFrameTff7b72> Tff7b72= Te6edf3_receivedDataTb4b4b4.Te6edf3receiveAsFlowTb4b4b4(Tb4b4b4)
T8b949e/** Running count of frames dropped because the queue was full. Diagnostic only; see [enqueueReceivedData]. */
Tff7b72private Tff7b72var Te6edf3droppedFrameCount Tff7b72= T79c0ff0L
Tff7b72private Tff7b72val Te6edf3_meshActivity Tff7b72=
Te6edf3MutableSharedFlowTff7b72<Te6edf3MeshActivityTff7b72>Tb4b4b4(Te6edf3extraBufferCapacity Tff7b72= T79c0ff6T79c0ff4Tb4b4b4, Te6edf3onBufferOverflow Tff7b72= Te6edf3BufferOverflowTb4b4b4.Te6edf3DROP_OLDESTTb4b4b4)
Tff7b72override Tff7b72val Te6edf3meshActivityTb4b4b4: Te6edf3FlowTff7b72<Te6edf3MeshActivityTff7b72> Tff7b72= Te6edf3_meshActivityTb4b4b4.Te6edf3asFlowTb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3_connectionError Tff7b72= Te6edf3MutableSharedFlowTff7b72<Tffa657StringTff7b72>Tb4b4b4(Te6edf3extraBufferCapacity Tff7b72= T79c0ff6T79c0ff4Tb4b4b4)
Tff7b72override Tff7b72val Te6edf3connectionErrorTb4b4b4: Te6edf3FlowTff7b72<Tffa657StringTff7b72> Tff7b72= Te6edf3_connectionErrorTb4b4b4.Te6edf3asFlowTb4b4b4(Tb4b4b4)
Tff7b72override Tff7b72val Te6edf3serviceScopeTb4b4b4: Te6edf3CoroutineScope
Tff7b72getTb4b4b4(Tb4b4b4) Tff7b72= Te6edf3_serviceScope
Tff7b72private Tff7b72var Te6edf3_serviceScope Tff7b72= Te6edf3CoroutineScopeTb4b4b4(Te6edf3dispatchersTb4b4b4.Te6edf3io Tff7b72+ Te6edf3SupervisorJobTb4b4b4(Tb4b4b4)Tb4b4b4)
Tff7b72private Tff7b72var Te6edf3radioTransportTb4b4b4: Te6edf3RadioTransport? Tff7b72= Tff7b72null
Tff7b72private Tff7b72var Te6edf3runningTransportIdTb4b4b4: Te6edf3InterfaceId? Tff7b72= Tff7b72null
Tff7b72private Tff7b72var Te6edf3isStarted Tff7b72= Tff7b72false
T8b949e/**
* Set while [stopTransportLocked] is draining the polite disconnect frame. [sendToRadio] checks this so any late
* traffic submitted after we've announced disconnection is dropped rather than racing in front of the firmware-side
* link teardown.
*/
Tf0883e@Volatile Tff7b72private Tff7b72var Te6edf3isStopping Tff7b72= Tff7b72false
T8b949e/**
* True while an explicit connection lifecycle is active (set by [connect]/[setDeviceAddress], cleared by
* [disconnect]). The hardware ([bluetoothRepository.state]) and network ([networkRepository.networkAvailable])
* listeners and the [checkLiveness] zombie-recovery path consult this to avoid starting a transport the user has
* torn down β without it, BT/network recovery emissions can wake a transport after explicit disconnect, leaving the
* app "connected" with no orchestrator collector and an unloaded NodeDB/channels.
*
* Guarded by [transportMutex]; every read/write site holds the lock. The @Volatile keeps diagnostic reads honest.
*/
Tf0883e@Volatile Tff7b72private Tff7b72var Te6edf3connectionRequested Tff7b72= Tff7b72false
Tff7b72private Tff7b72val Te6edf3gattCacheInvalidationRequested Tff7b72= Te6edf3atomicTb4b4b4(Tff7b72falseTb4b4b4)
T8b949e/** Prevents concurrent liveness-induced transport restarts from stacking. */
Tff7b72private Tff7b72val Te6edf3isRestarting Tff7b72= Te6edf3atomicTb4b4b4(Tff7b72falseTb4b4b4)
Tff7b72private Tff7b72val Te6edf3listenersInitialized Tff7b72= Te6edf3atomicTb4b4b4(Tff7b72falseTb4b4b4)
Tff7b72private Tff7b72var Te6edf3heartbeatJobTb4b4b4: Te6edf3Job? Tff7b72= Tff7b72null
Tff7b72private Tff7b72var Te6edf3lastHeartbeatMillis Tff7b72= T79c0ff0L
Tf0883e@Volatile Tff7b72private Tff7b72var Te6edf3lastDataReceivedMillis Tff7b72= T79c0ff0L
T8b949e/**
* Internal test seam for deterministic clock injection. Production uses [nowMillis]; tests override this to a
* controllable clock so [onConnect], [handleFromRadio], [checkLiveness], and [keepAlive] all share one coherent
* time source. Not a constructor parameter to avoid breaking Koin @Single annotation generation (which would try to
* resolve `() -> Long` from the DI graph).
*/
Tf0883e@Volatile
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffMemberVisibilityCanBePrivateTa5d6ff"Tb4b4b4)
Tff7b72internal Tff7b72var Te6edf3clockMillisTb4b4b4: Tb4b4b4(Tb4b4b4) Tff7b72-Tff7b72> Tffa657Long Tff7b72= Tb4b4b4{ Te6edf3nowMillis Tb4b4b4}
T8b949e/** The current time from the injected clock. */
Tff7b72private Tff7b72fun Td2a8ffnowTb4b4b4(Tb4b4b4)Tb4b4b4: Tffa657Long Tff7b72= Te6edf3clockMillisTb4b4b4(Tb4b4b4)
Tff7b72companion Tff7b72object Tb4b4b4{
T8b949e/**
* Capacity of the inbound frame queue.
*
* Sized against the burst a connect produces, which is roughly `35 + 5N` frames for a NodeDB of N nodes:
* - ~35 fixed: `my_info`, metadata, ~10 config, ~14 moduleConfig, 8 channels, deviceui, `config_complete`.
* - N thin `NodeInfo` frames β one per *hot-store* node. Warm-tier entries (firmware `WARM_NODE_COUNT`, up to
* 2000 on a native host) are identity-only records for evicted nodes and are NOT streamed to the phone, so
* they do not contribute here.
* - Up to 4N more from the post-`config_complete` replay drain, which re-sends stored position, telemetry,
* environment and status as ordinary mesh packets β one of each per node.
*
* N is firmware `MAX_NUM_NODES`: 250 on portduino/native-host and top-tier ESP32-S3, 120 on nRF52840 and
* generic ESP32, 10 on STM32WL. So the realistic worst case is a Linux/Pi node at N=250 β ~1285 frames, and
* those arrive over TCP, the transport most able to outrun the consumer. 8192 keeps roughly 6x headroom over
* that; a custom build raising `MAX_NUM_NODES` past ~1630 would need this raised too.
*
* Memory stays bounded at capacity x frame size. On the stream transports a frame cannot exceed
* `StreamFrameCodec.MAX_TO_FROM_RADIO_SIZE` (512 B); the BLE path passes through whatever the GATT read
* returned, which ATT caps at 512 B in practice rather than by anything enforced here. A thin `NodeInfo` is
* closer to 100 B, so the realistic ceiling is well under the ~4 MB absolute worst case.
*/
Tff7b72const Tff7b72val Te6edf3RECEIVE_QUEUE_CAPACITY Tff7b72= T79c0ff8T79c0ff1T79c0ff9T79c0ff2
T8b949e/** Log one dropped-frame warning per this many drops. See [enqueueReceivedData]. */
Tff7b72private Tff7b72const Tff7b72val Te6edf3DROP_LOG_INTERVAL Tff7b72= T79c0ff5T79c0ff1T79c0ff2L
T8b949e/**
* Per-frame ceiling, matching `StreamFrameCodec.MAX_TO_FROM_RADIO_SIZE`.
*
* Duplicated rather than imported because `core:service` does not depend on `core:network`; the stream codec
* enforces the same number on its own path, and ATT caps BLE at the same value in practice.
*/
Tff7b72const Tff7b72val Te6edf3MAX_FRAME_BYTES Tff7b72= T79c0ff5T79c0ff1T79c0ff2
Tff7b72private Tff7b72const Tff7b72val Te6edf3HEARTBEAT_INTERVAL_MILLIS Tff7b72= T79c0ff3T79c0ff0 Tff7b72* T79c0ff1T79c0ff0T79c0ff0T79c0ff0L
T8b949e// If we haven't received any data from the radio within this window after sending a
T8b949e// heartbeat while the connection is nominally "Connected", the connection is likely a
T8b949e// zombie (BLE stack didn't report disconnect). Two missed heartbeat intervals gives
T8b949e// the firmware a reasonable window to respond or send telemetry.
Tff7b72private Tff7b72const Tff7b72val Te6edf3LIVENESS_TIMEOUT_MILLIS Tff7b72= Te6edf3HEARTBEAT_INTERVAL_MILLIS Tff7b72* T79c0ff2
T8b949e/**
* Upper bound on how long we wait for the polite `ToRadio(disconnect = true)` frame to flush before tearing the
* transport down. 500ms gives BLE's write-retry path (`BleRetry` backs off 500ms) room for one attempt on a
* flaky GATT connection. Serial and TCP typically flush well under this window.
*/
Tff7b72private Tff7b72const Tff7b72val Te6edf3POLITE_DISCONNECT_DRAIN_MS Tff7b72= T79c0ff5T79c0ff0T79c0ff0L
Tb4b4b4}
Tff7b72private Tff7b72val Te6edf3initLock Tff7b72= Te6edf3MutexTb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3transportMutex Tff7b72= Te6edf3MutexTb4b4b4(Tb4b4b4)
Tff7b72init Tb4b4b4{
T8b949e// Address sync runs INDEPENDENTLY of connect() so that observers (notably
T8b949e// MeshServiceOrchestrator) see a valid currentDeviceAddressFlow BEFORE the transport
T8b949e// starts. Previously this mirror lived inside initStateListeners()'s devAddr listener and
T8b949e// was coupled to startTransportLocked() β but initStateListeners() is only invoked from
T8b949e// connect(), so the flow never updated until connect() ran. That created a cold-start race:
T8b949e// the orchestrator's currentDeviceAddressFlow observer could fire AFTER the transport had
T8b949e// already started, violating the invariant that the active DB must be switched to the
T8b949e// selected device's DB before its transport starts.
T8b949e//
T8b949e// This listener ONLY mirrors radioPrefs.devAddr into _currentDeviceAddressFlow; it never
T8b949e// starts a transport. Transport start remains driven exclusively by connect() (initial),
T8b949e// setDeviceAddress() (explicit user switch), BLE/network state changes (environment
T8b949e// recovery), and liveness restarts (zombie recovery) β see startTransportLocked() callers.
T8b949e// _currentDeviceAddressFlow is a MutableStateFlow (atomic .value), so the unconditional
T8b949e// assignment here is race-free without holding transportMutex; same-address writes are
T8b949e// idempotent no-ops.
Te6edf3radioPrefsTb4b4b4.Te6edf3devAddr
Tb4b4b4.Te6edf3onEach Tb4b4b4{ Te6edf3addr Tff7b72-Tff7b72> Te6edf3_currentDeviceAddressFlowTb4b4b4.Te6edf3value Tff7b72= Te6edf3addr Tb4b4b4}
Tb4b4b4.Te6edf3catch Tb4b4b4{ Te6edf3LoggerTb4b4b4.Te6edf3eTb4b4b4(Tffa657itTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffradioPrefs.devAddr address-sync flow crashedTa5d6ff" Tb4b4b4} Tb4b4b4}
Tb4b4b4.Te6edf3launchInTb4b4b4(Te6edf3processLifecycleTb4b4b4.Te6edf3coroutineScopeTb4b4b4)
Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffinitStateListenersTb4b4b4(Tb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3listenersInitializedTb4b4b4.Te6edf3valueTb4b4b4) Tff7b72return
Te6edf3processLifecycleTb4b4b4.Te6edf3coroutineScopeTb4b4b4.Te6edf3launch Tb4b4b4{
Te6edf3initLockTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3listenersInitializedTb4b4b4.Te6edf3valueTb4b4b4) Tff7b72returnTf0883e@withLock
Te6edf3listenersInitializedTb4b4b4.Te6edf3value Tff7b72= Tff7b72true
Te6edf3bluetoothRepositoryTb4b4b4.Te6edf3state
Tb4b4b4.Te6edf3onEach Tb4b4b4{ Te6edf3state Tff7b72-Tff7b72>
Te6edf3transportMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3stateTb4b4b4.Te6edf3enabledTb4b4b4) Tb4b4b4{
T8b949e// Environmental recovery only: don't wake a transport the user has
T8b949e// explicitly disconnected from. stopTransportLocked() below still fires on
T8b949e// BLE-disabled to tear down a running BLE link, but we deliberately do NOT
T8b949e// clear connectionRequested here β that is disconnect()'s job.
Tff7b72if Tb4b4b4(Te6edf3connectionRequestedTb4b4b4) Tb4b4b4{
Te6edf3startTransportLockedTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tb4b4b4} Tff7b72else Tff7b72if Tb4b4b4(Te6edf3runningTransportId Tff7b72=Tff7b72= Te6edf3InterfaceIdTb4b4b4.Te6edf3BLUETOOTHTb4b4b4) Tb4b4b4{
Te6edf3stopTransportLockedTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4.Te6edf3catch Tb4b4b4{ Te6edf3LoggerTb4b4b4.Te6edf3eTb4b4b4(Tffa657itTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffbluetoothRepository.state flow crashedTa5d6ff" Tb4b4b4} Tb4b4b4}
Tb4b4b4.Te6edf3launchInTb4b4b4(Te6edf3processLifecycleTb4b4b4.Te6edf3coroutineScopeTb4b4b4)
Te6edf3networkRepositoryTb4b4b4.Te6edf3networkAvailable
Tb4b4b4.Te6edf3onEach Tb4b4b4{ Te6edf3state Tff7b72-Tff7b72>
Te6edf3transportMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3stateTb4b4b4) Tb4b4b4{
T8b949e// Environmental recovery only β see the BLE listener above for rationale.
Tff7b72if Tb4b4b4(Te6edf3connectionRequestedTb4b4b4) Tb4b4b4{
Te6edf3startTransportLockedTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tb4b4b4} Tff7b72else Tff7b72if Tb4b4b4(Te6edf3runningTransportId Tff7b72=Tff7b72= Te6edf3InterfaceIdTb4b4b4.Te6edf3TCPTb4b4b4) Tb4b4b4{
Te6edf3stopTransportLockedTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4.Te6edf3catch Tb4b4b4{ Te6edf3LoggerTb4b4b4.Te6edf3eTb4b4b4(Tffa657itTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffnetworkRepository.networkAvailable flow crashedTa5d6ff" Tb4b4b4} Tb4b4b4}
Tb4b4b4.Te6edf3launchInTb4b4b4(Te6edf3processLifecycleTb4b4b4.Te6edf3coroutineScopeTb4b4b4)
Te6edf3observeUsbRecoveryTriggersTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Observes selected-SERIAL-device presence and transport state to auto-restart USB serial after replug.
*
* Recovery arms only after the selected serial key is observed absent, then fires when that same key is present
* while the transport is in [ConnectionState.DeviceSleep]. This prevents normal unplug races from restarting
* against stale presence before UsbRepository removes the key.
*
* Including [_connectionState] closes the race where presence returns before the I/O-death callback flips state to
* [ConnectionState.DeviceSleep]. Each physical replug consumes one armed edge, and recovery remains gated by
* [connectionRequested] and [runningTransportId] so explicit disconnects cannot resurrect a transport.
*/
Tff7b72private Tff7b72fun Td2a8ffobserveUsbRecoveryTriggersTb4b4b4(Tb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3selectedPresence Tff7b72=
Te6edf3combineTb4b4b4(Te6edf3currentDeviceAddressFlowTb4b4b4, Te6edf3serialDevicePresenceTb4b4b4.Te6edf3deviceKeysTb4b4b4, Tff7b72::Te6edf3selectedSerialPresenceTb4b4b4)
Tb4b4b4.Te6edf3distinctUntilChangedTb4b4b4(Tb4b4b4)
Te6edf3combineTb4b4b4(Te6edf3selectedPresenceTb4b4b4, Te6edf3_connectionStateTb4b4b4, Tff7b72::Te6edf3UsbRecoverySnapshotTb4b4b4)
Tb4b4b4.Te6edf3runningFoldTb4b4b4(Te6edf3UsbRecoveryTriggerStateTb4b4b4(Tb4b4b4)Tb4b4b4) Tb4b4b4{ Te6edf3triggerStateTb4b4b4, Te6edf3snapshot Tff7b72-Tff7b72> Te6edf3triggerStateTb4b4b4.Te6edf3nextTb4b4b4(Te6edf3snapshotTb4b4b4) Tb4b4b4}
Tb4b4b4.Te6edf3distinctUntilChangedTb4b4b4(Tb4b4b4)
Tb4b4b4.Te6edf3dropTb4b4b4(T79c0ff1Tb4b4b4)
Tb4b4b4.Te6edf3filter Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3triggerRecovery Tb4b4b4}
Tb4b4b4.Te6edf3onEach Tb4b4b4{
Te6edf3transportMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3connectionRequestedTb4b4b4) Tff7b72returnTf0883e@withLock
Tff7b72if Tb4b4b4(Te6edf3runningTransportId Tff7b72!Tff7b72= Te6edf3InterfaceIdTb4b4b4.Te6edf3SERIALTb4b4b4) Tff7b72returnTf0883e@withLock
T8b949e// Race-defense: the combine snapshot may be stale by the time we acquire
T8b949e// transportMutex β another path (setDeviceAddress, BLE liveness restart) may
T8b949e// have just brought this transport up to Connected/Connecting. The combine-level
T8b949e// state filter narrows the trigger; this check guards the emission β mutex
T8b949e// acquisition window so we never tear down a fresh healthy transport.
Tff7b72val Te6edf3state Tff7b72= Te6edf3_connectionStateTb4b4b4.Te6edf3value
Tff7b72if Tb4b4b4(Te6edf3state Tff7b72is Te6edf3ConnectionStateTb4b4b4.Te6edf3Connected Tff7b72|Tff7b72| Te6edf3state Tff7b72is Te6edf3ConnectionStateTb4b4b4.Te6edf3ConnectingTb4b4b4) Tb4b4b4{
Tff7b72returnTf0883e@withLock
Tb4b4b4}
T8b949e// Re-check presence under the lock β the combine snapshot may be stale
T8b949e// if the device was unplugged or the selection changed while awaiting
T8b949e// the mutex. Mirrors the race-defense pattern of the state check above.
Tff7b72val Te6edf3currentKeys Tff7b72= Te6edf3serialDevicePresenceTb4b4b4.Te6edf3deviceKeysTb4b4b4.Te6edf3value
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3selectedSerialPresenceTb4b4b4(Te6edf3_currentDeviceAddressFlowTb4b4b4.Te6edf3valueTb4b4b4, Te6edf3currentKeysTb4b4b4)Tb4b4b4.Te6edf3presentTb4b4b4) Tb4b4b4{
Tff7b72returnTf0883e@withLock
Tb4b4b4}
T8b949e// A previously-started SERIAL transport that died on unplug is still held in
T8b949e// radioTransport (the I/O-death path emits DeviceSleep but does not null it). Tear it
T8b949e// down silently so startTransportLocked() can build a fresh one for the replugged
T8b949e// device. Mirrors the BLE liveness recovery shape: notifyPermanent=false (no
T8b949e// user-facing Disconnected), sendPoliteDisconnect=false (the link is already gone).
T8b949e// Keep teardown and bring-up isolated: a transient USB close error must not skip
T8b949e// the fresh start, and neither error should terminate this long-lived recovery flow.
Te6edf3ignoreExceptionSuspend Tb4b4b4{
Te6edf3stopTransportLockedTb4b4b4(Te6edf3notifyPermanent Tff7b72= Tff7b72falseTb4b4b4, Te6edf3sendPoliteDisconnect Tff7b72= Tff7b72falseTb4b4b4)
Tb4b4b4}
Te6edf3ignoreExceptionSuspend Tb4b4b4{ Te6edf3startTransportLockedTb4b4b4(Tb4b4b4) Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4.Te6edf3catch Tb4b4b4{ Te6edf3LoggerTb4b4b4.Te6edf3eTb4b4b4(Tffa657itTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffserialDevicePresence recovery flow crashedTa5d6ff" Tb4b4b4} Tb4b4b4}
Tb4b4b4.Te6edf3launchInTb4b4b4(Te6edf3processLifecycleTb4b4b4.Te6edf3coroutineScopeTb4b4b4)
Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffconnectTb4b4b4(Tb4b4b4) Tb4b4b4{
Te6edf3processLifecycleTb4b4b4.Te6edf3coroutineScopeTb4b4b4.Te6edf3launch Tb4b4b4{
Te6edf3transportMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
T8b949e// Mark the connection lifecycle as active BEFORE starting so concurrent
T8b949e// hardware/network listeners observe the gate as open.
Te6edf3connectionRequested Tff7b72= Tff7b72true
T8b949e// connect() is fire-and-forget, so a recoverable factory failure has no caller to receive it. The
T8b949e// start path already rolls back partial session state; contain the failure here so it does not become
T8b949e// an uncaught lifecycle-scope exception, while CancellationException and fatal Errors still propagate.
Te6edf3ignoreExceptionSuspend Tb4b4b4{ Te6edf3startTransportLockedTb4b4b4(Tb4b4b4) Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Te6edf3initStateListenersTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tff7b72override Tff7b72suspend Tff7b72fun Td2a8ffdisconnectTb4b4b4(Tb4b4b4) Tb4b4b4{
Te6edf3transportMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
T8b949e// Tear the gate down BEFORE stopTransportLocked() so a concurrent state-listener
T8b949e// emission arriving while we wait for the mutex cannot re-start the transport.
Te6edf3connectionRequested Tff7b72= Tff7b72false
Te6edf3gattCacheInvalidationRequestedTb4b4b4.Te6edf3value Tff7b72= Tff7b72false
Te6edf3ignoreExceptionSuspend Tb4b4b4{ Te6edf3stopTransportLockedTb4b4b4(Tb4b4b4) Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72suspend Tff7b72fun Td2a8ffrestartTransportTb4b4b4(Tb4b4b4) Tb4b4b4{
T8b949e// CAS BEFORE the mutex, mirroring checkLiveness()'s coordination structure: both
T8b949e// restart paths CAS synchronously, one wins, one loses immediately. Performing the
T8b949e// CAS inside transportMutex.withLock races checkLiveness's outer
T8b949e// `finally { isRestarting = false }` (which runs AFTER mutex release): a queued
T8b949e// restartTransport that resumes from mutex.wait can observe isRestarting == false,
T8b949e// win the CAS, and produce an extra transport cycle (3 instead of 2) under the JVM's
T8b949e// real dispatcher. The loser here observes isRestarting == true and defers to the
T8b949e// in-flight cycle. startTransportLocked() is idempotent w.r.t. an existing transport,
T8b949e// but the CAS also prevents a double stop/stop race on the teardown side.
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3isRestartingTb4b4b4.Te6edf3compareAndSetTb4b4b4(Te6edf3expect Tff7b72= Tff7b72falseTb4b4b4, Te6edf3update Tff7b72= Tff7b72trueTb4b4b4)Tb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffrestartTransport: skipped, concurrent restart in progressTa5d6ff" Tb4b4b4}
Tff7b72return
Tb4b4b4}
Tff7b72try Tb4b4b4{
Te6edf3transportMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
T8b949e// Silent recovery for app-level handshake stalls. The transport may still be physically
T8b949e// up (TCP socket alive, firmware unresponsive to want_config_id), so cycling it in place
T8b949e// WITHOUT clearing connectionRequested avoids the split-brain where setDeviceAddress's
T8b949e// fast-path would otherwise block same-node reconnect. The caller (MeshConnectionManager)
T8b949e// is responsible for the app-level Disconnected flip; this method only cycles the
T8b949e// transport and emits the transport-level DeviceSleep -> Connected transitions via
T8b949e// callbacks (no Connecting β that is an app-level state, not a transport emission).
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3connectionRequestedTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffrestartTransport: skipped-not-requestedTa5d6ff" Tb4b4b4}
Tff7b72returnTf0883e@withLock
Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3getBondedDeviceAddressTb4b4b4(Tb4b4b4) Tff7b72=Tff7b72= Tff7b72nullTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffrestartTransport: skipped-no-addressTa5d6ff" Tb4b4b4}
Tff7b72returnTf0883e@withLock
Tb4b4b4}
T8b949e// Honor the documented "safe no-op when no transport running" contract: environmental
T8b949e// stops (network unavailable, BLE disabled) intentionally preserve
T8b949e// connectionRequested=true so the recovery listeners above can re-bring-up the
T8b949e// transport later. A stale restart job running after such a stop must NOT bypass that
T8b949e// recovery path by creating a transport directly via startTransportLocked().
Tff7b72if Tb4b4b4(Te6edf3radioTransport Tff7b72=Tff7b72= Tff7b72nullTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffrestartTransport: skipped-no-transportTa5d6ff" Tb4b4b4}
Tff7b72returnTf0883e@withLock
Tb4b4b4}
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffrestartTransport: restarting transport for Tffd700${Te6edf3getDeviceAddressTb4b4b4(Tb4b4b4)Tff7b72?.Te6edf3anonymizeTffd700}Ta5d6ff" Tb4b4b4}
T8b949e// Mirror checkLiveness()'s coordination contract: emit a transport-level
T8b949e// Connected -> DeviceSleep transition before the stop/start cycle so
T8b949e// transport-level observers see a full DeviceSleep -> Connected cycle when
T8b949e// the fresh transport's onConnect() fires and can re-trigger their onConnected
T8b949e// logic. The caller (runSiblingHandshakeRecovery) flips the app-level state to
T8b949e// Disconnected before invoking restartTransport(), so this emission exists for
T8b949e// transport-level coordination, not for app-level StateFlow dedupe. Silent:
T8b949e// transient DeviceSleep, not a permanent Disconnected.
Te6edf3onDisconnectTb4b4b4(Te6edf3isPermanent Tff7b72= Tff7b72falseTb4b4b4)
T8b949e// notifyPermanent=false below (no user-facing Disconnected modal β the app-level
T8b949e// state machine drives that separately) and sendPoliteDisconnect=false (firmware
T8b949e// is unresponsive, writing a goodbye frame into a dead link only delays teardown).
Te6edf3ignoreExceptionSuspend Tb4b4b4{ Te6edf3stopTransportLockedTb4b4b4(Te6edf3notifyPermanent Tff7b72= Tff7b72falseTb4b4b4, Te6edf3sendPoliteDisconnect Tff7b72= Tff7b72falseTb4b4b4) Tb4b4b4}
T8b949e// Defense-in-depth mirroring the liveness recovery gate; today all
T8b949e// connectionRequested mutators (connect/disconnect/setDeviceAddress) hold
T8b949e// transportMutex so this re-check is unreachable, but it guards against future
T8b949e// refactors that mutate the gate without serialization.
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3connectionRequestedTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffrestartTransport: aborted, disconnect requested during stopTa5d6ff" Tb4b4b4}
Tff7b72returnTf0883e@withLock
Tb4b4b4}
T8b949e// Drop whatever the dead session left queued before admitting the replacement. The consumer discards
T8b949e// stale-generation frames on dequeue, but they still occupy slots until then β and now that the queue
T8b949e// is bounded, a backlog carried across the cycle can make the fresh session's handshake frames fail
T8b949e// trySend. `MeshServiceOrchestrator.start()` drains for its own stop/start path, but a transport-level
T8b949e// restart does not go through it, so the drain has to happen here too.
Te6edf3resetReceivedBufferTb4b4b4(Tb4b4b4)
T8b949e// startTransportLocked() re-validates the selected address (no-op if null) and emits
T8b949e// Connected through the transport callbacks (via the new transport's onConnect) once
T8b949e// the fresh transport comes up β there is no Connecting emission at the transport
T8b949e// layer (that is an app-level state owned by MeshConnectionManager).
Te6edf3startTransportLockedTb4b4b4(Tb4b4b4)
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{ Ta5d6ff"Ta5d6ffrestartTransport: completedTa5d6ff" Tb4b4b4}
Tb4b4b4}
Tb4b4b4} Tff7b72finally Tb4b4b4{
Te6edf3isRestartingTb4b4b4.Te6edf3value Tff7b72= Tff7b72false
Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffrequestGattCacheInvalidationOnNextConnectTb4b4b4(Tb4b4b4) Tb4b4b4{
Te6edf3gattCacheInvalidationRequestedTb4b4b4.Te6edf3value Tff7b72= Tff7b72true
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffGATT cache invalidation requested for next BLE connectTa5d6ff" Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffconsumeGattCacheInvalidationRequestTb4b4b4(Tb4b4b4)Tb4b4b4: Tffa657Boolean Tff7b72= Te6edf3gattCacheInvalidationRequestedTb4b4b4.Te6edf3getAndSetTb4b4b4(Tff7b72falseTb4b4b4)
Tff7b72override Tff7b72fun Td2a8ffisMockTransportTb4b4b4(Tb4b4b4)Tb4b4b4: Tffa657Boolean Tff7b72= Te6edf3transportFactoryTb4b4b4.Te6edf3isMockTransportTb4b4b4(Tb4b4b4)
Tff7b72override Tff7b72fun Td2a8fftoInterfaceAddressTb4b4b4(Te6edf3interfaceIdTb4b4b4: Te6edf3InterfaceIdTb4b4b4, Te6edf3restTb4b4b4: Tffa657StringTb4b4b4)Tb4b4b4: Tffa657String Tff7b72=
Te6edf3transportFactoryTb4b4b4.Te6edf3toInterfaceAddressTb4b4b4(Te6edf3interfaceIdTb4b4b4, Te6edf3restTb4b4b4)
Tff7b72override Tff7b72fun Td2a8ffgetDeviceAddressTb4b4b4(Tb4b4b4)Tb4b4b4: Tffa657String? Tff7b72= Te6edf3_currentDeviceAddressFlowTb4b4b4.Te6edf3value
Tff7b72private Tff7b72fun Td2a8ffgetBondedDeviceAddressTb4b4b4(Tb4b4b4)Tb4b4b4: Tffa657String? Tb4b4b4{
Tff7b72val Te6edf3address Tff7b72= Te6edf3getDeviceAddressTb4b4b4(Tb4b4b4)
Tff7b72return Tff7b72if Tb4b4b4(Te6edf3transportFactoryTb4b4b4.Te6edf3isAddressValidTb4b4b4(Te6edf3addressTb4b4b4)Tb4b4b4) Tb4b4b4{
Te6edf3address
Tb4b4b4} Tff7b72else Tb4b4b4{
Tff7b72null
Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffsetDeviceAddressTb4b4b4(Te6edf3deviceAddrTb4b4b4: Tffa657String?Tb4b4b4)Tb4b4b4: Tffa657Boolean Tb4b4b4{
Tff7b72val Te6edf3sanitized Tff7b72= Tff7b72if Tb4b4b4(Te6edf3deviceAddr Tff7b72=Tff7b72= Ta5d6ff"Ta5d6ffnTa5d6ff" Tff7b72|Tff7b72| Te6edf3deviceAddrTb4b4b4.Te6edf3isNullOrBlankTb4b4b4(Tb4b4b4)Tb4b4b4) Tff7b72null Tff7b72else Te6edf3deviceAddr
Tff7b72if Tb4b4b4(Te6edf3getBondedDeviceAddressTb4b4b4(Tb4b4b4) Tff7b72=Tff7b72= Te6edf3sanitized Tff7b72&Tff7b72& Te6edf3isStarted Tff7b72&Tff7b72& Te6edf3_connectionStateTb4b4b4.Te6edf3value Tff7b72=Tff7b72= Te6edf3ConnectionStateTb4b4b4.Te6edf3ConnectedTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffIgnoring setBondedDevice Tffd700${Te6edf3sanitizedTff7b72?.Te6edf3anonymizeTffd700}Ta5d6ff, already using that deviceTa5d6ff" Tb4b4b4}
Tff7b72return Tff7b72false
Tb4b4b4}
Tff7b72val Te6edf3previousAddress Tff7b72= Te6edf3getBondedDeviceAddressTb4b4b4(Tb4b4b4)
Te6edf3analyticsTb4b4b4.Te6edf3trackTb4b4b4(Ta5d6ff"Ta5d6ffmesh_bondTa5d6ff"Tb4b4b4)
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffSetting bonded device to Tffd700${Te6edf3sanitizedTff7b72?.Te6edf3anonymizeTffd700}Ta5d6ff" Tb4b4b4}
Te6edf3radioPrefsTb4b4b4.Te6edf3setDevAddrTb4b4b4(Te6edf3sanitizedTb4b4b4)
Te6edf3_currentDeviceAddressFlowTb4b4b4.Te6edf3value Tff7b72= Te6edf3sanitized
Te6edf3processLifecycleTb4b4b4.Te6edf3coroutineScopeTb4b4b4.Te6edf3launch Tb4b4b4{
Te6edf3transportMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
T8b949e// The sanitized address is the single source of truth for the connectionRequested
T8b949e// gate: a real address arms the lifecycle (connect() equivalent) so environmental
T8b949e// listeners cannot race the rebind into a "down" state; null/("n") is a deselect
T8b949e// that MUST clear the gate so subsequent BT/network recovery cannot resurrect a
T8b949e// transport for a device the user explicitly tore down. Only start a fresh
T8b949e// transport when an address was actually selected.
Te6edf3connectionRequested Tff7b72= Te6edf3sanitized Tff7b72!Tff7b72= Tff7b72null
Tff7b72if Tb4b4b4(Te6edf3sanitized Tff7b72!Tff7b72= Tff7b72null Tff7b72&Tff7b72& Te6edf3previousAddress Tff7b72!Tff7b72= Tff7b72null Tff7b72&Tff7b72& Te6edf3sanitized Tff7b72!Tff7b72= Te6edf3previousAddressTb4b4b4) Tb4b4b4{
Te6edf3gattCacheInvalidationRequestedTb4b4b4.Te6edf3value Tff7b72= Tff7b72false
Tb4b4b4}
Te6edf3ignoreExceptionSuspend Tb4b4b4{ Te6edf3stopTransportLockedTb4b4b4(Tb4b4b4) Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3sanitized Tff7b72!Tff7b72= Tff7b72nullTb4b4b4) Tb4b4b4{
T8b949e// setDeviceAddress() is fire-and-forget. startTransportLocked() has already rolled back any
T8b949e// partially admitted session, so contain a recoverable factory failure instead of crashing the
T8b949e// process-lifecycle scope. Explicit suspend restart callers still receive their failures.
Te6edf3ignoreExceptionSuspend Tb4b4b4{ Te6edf3startTransportLockedTb4b4b4(Tb4b4b4) Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72return Tff7b72true
Tb4b4b4}
T8b949e/** Must be called under [transportMutex]. */
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4)
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffstartTransportLockedTb4b4b4(Tb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3radioTransport Tff7b72!Tff7b72= Tff7b72nullTb4b4b4) Tff7b72return
T8b949e// Never autoconnect to the simulated node. The mock transport may be offered in the
T8b949e// device-picker UI on debug builds, but it must only connect when the user explicitly
T8b949e// selects it (i.e. its address is stored in radioPrefs).
Tff7b72val Te6edf3address Tff7b72= Te6edf3getBondedDeviceAddressTb4b4b4(Tb4b4b4)
Tff7b72if Tb4b4b4(Te6edf3address Tff7b72=Tff7b72= Tff7b72nullTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffNo valid address to connect toTa5d6ff" Tb4b4b4}
Tff7b72return
Tb4b4b4}
T8b949e// Build a fresh per-instance session and admit it BEFORE constructing the transport so the wrapper captures
T8b949e// the new session. Bumping the public generation also signals downstream consumers (e.g.
T8b949e// RadioControllerImpl)
T8b949e// to invalidate session-scoped state from any previous transport instance.
Tff7b72val Te6edf3generation Tff7b72= Te6edf3sessionGenerationCounterTb4b4b4.Te6edf3incrementAndGetTb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3session Tff7b72= Te6edf3RadioTransportSessionTb4b4b4(Te6edf3generation Tff7b72= Te6edf3generationTb4b4b4, Te6edf3address Tff7b72= Te6edf3addressTb4b4b4)
Te6edf3synchronizedTb4b4b4(Te6edf3sessionCallbackLockTb4b4b4) Tb4b4b4{
Te6edf3checkTb4b4b4(Te6edf3activeTransportSession Tff7b72=Tff7b72= Tff7b72null Tff7b72&Tff7b72& Te6edf3admittedSessionOperations Tff7b72=Tff7b72= T79c0ff0Tb4b4b4) Tb4b4b4{
Ta5d6ff"Ta5d6ffCannot admit a transport while the previous session is still drainingTa5d6ff"
Tb4b4b4}
Te6edf3activeTransportSession Tff7b72= Te6edf3session
Te6edf3sessionAdmissionOpen Tff7b72= Tff7b72true
Te6edf3_activeSessionTb4b4b4.Te6edf3value Tff7b72= Te6edf3sessionTb4b4b4.Te6edf3context
Te6edf3_sessionGenerationTb4b4b4.Te6edf3value Tff7b72= Te6edf3generation
Tb4b4b4}
Tff7b72val Te6edf3sessionBoundService Tff7b72= Te6edf3SessionBoundRadioInterfaceServiceTb4b4b4(Te6edf3sessionTb4b4b4)
Tff7b72val Te6edf3connectionStateBeforeStart Tff7b72= Te6edf3_connectionStateTb4b4b4.Te6edf3value
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{ Ta5d6ff"Ta5d6ffStarting radio transport for Tffd700${Te6edf3addressTb4b4b4.Te6edf3anonymizeTffd700}Ta5d6ff (generation=Tffd700$Te6edf3generationTa5d6ff)Ta5d6ff" Tb4b4b4}
Tff7b72val Te6edf3newTransport Tff7b72=
Tff7b72try Tb4b4b4{
Te6edf3transportFactoryTb4b4b4.Te6edf3createTransportTb4b4b4(Te6edf3addressTb4b4b4, Te6edf3sessionBoundServiceTb4b4b4)
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3failureTb4b4b4: Te6edf3ThrowableTb4b4b4) Tb4b4b4{
T8b949e// A failed factory call consumed a generation that may already have reached observers. Keep the
T8b949e// generation monotonic, but revoke the admitted session and every transport-lifecycle field before
T8b949e// rethrowing the original failure. A later retry will receive a strictly newer generation.
Te6edf3revokeTransportSessionTb4b4b4(Te6edf3sessionTb4b4b4)
Te6edf3radioTransport Tff7b72= Tff7b72null
Te6edf3runningTransportId Tff7b72= Tff7b72null
Te6edf3isStarted Tff7b72= Tff7b72false
Te6edf3_connectionStateTb4b4b4.Te6edf3value Tff7b72= Te6edf3connectionStateBeforeStart
Tff7b72throw Te6edf3failure
Tb4b4b4}
Te6edf3radioTransport Tff7b72= Te6edf3newTransport
Te6edf3runningTransportId Tff7b72= Te6edf3addressTb4b4b4.Te6edf3firstOrNullTb4b4b4(Tb4b4b4)Tff7b72?.Te6edf3let Tb4b4b4{ Te6edf3InterfaceIdTb4b4b4.Te6edf3forIdCharTb4b4b4(Tffa657itTb4b4b4) Tb4b4b4}
Te6edf3isStarted Tff7b72= Tff7b72true
Te6edf3startHeartbeatTb4b4b4(Tb4b4b4)
Tb4b4b4}
T8b949e/**
* Must be called under [transportMutex].
*
* @param notifyPermanent When `true`, emits a permanent disconnect state to [connectionState]. Set `false` during
* automatic liveness recovery to avoid surfacing a user-facing disconnect.
* @param sendPoliteDisconnect When `true`, sends a `ToRadio(disconnect = true)` frame to the firmware before
* tearing down. Set `false` when the transport is already dead (zombie session) to avoid writing into a broken
* link.
*/
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffstopTransportLockedTb4b4b4(Te6edf3notifyPermanentTb4b4b4: Tffa657Boolean Tff7b72= Tff7b72trueTb4b4b4, Te6edf3sendPoliteDisconnectTb4b4b4: Tffa657Boolean Tff7b72= Tff7b72trueTb4b4b4) Tff7b72=
Te6edf3withContextTb4b4b4(Te6edf3NonCancellableTb4b4b4) Tb4b4b4{ Te6edf3finishTransportTeardownTb4b4b4(Te6edf3notifyPermanentTb4b4b4, Te6edf3sendPoliteDisconnectTb4b4b4) Tb4b4b4}
T8b949e/** Completes teardown after admission closes, even if its requester is cancelled while leases drain. */
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8fffinishTransportTeardownTb4b4b4(Te6edf3notifyPermanentTb4b4b4: Tffa657BooleanTb4b4b4, Te6edf3sendPoliteDisconnectTb4b4b4: Tffa657BooleanTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3currentTransport Tff7b72= Te6edf3radioTransport
Tff7b72val Te6edf3currentSession Tff7b72= Te6edf3synchronizedTb4b4b4(Te6edf3sessionCallbackLockTb4b4b4) Tb4b4b4{ Te6edf3activeTransportSession Tb4b4b4}
T8b949e// Reject queued callbacks and new suspend work immediately, then drain existing leases before admitting a
T8b949e// replacement generation or closing the old transport.
Te6edf3revokeTransportSessionTb4b4b4(Te6edf3currentSessionTb4b4b4)
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{ Ta5d6ff"Ta5d6ffStopping transport Tffd700$Te6edf3currentTransportTa5d6ff" Tb4b4b4}
T8b949e// Best-effort polite goodbye: tell the firmware we're disconnecting on purpose so it can
T8b949e// tear down its side of the link cleanly instead of relying on timeouts / hardware events.
T8b949e// Flip isStopping before sending so any concurrent sendToRadio() drops incoming traffic β
T8b949e// we don't want normal packets racing behind the disconnect frame. Skip only when already
T8b949e// Disconnected; firmware can still consume the goodbye while handshaking or sleeping, so
T8b949e// it's worth sending in every other state. The send is fire-and-forget through the
T8b949e// transport's own scope; the drain delay gives async transports a window to flush before
T8b949e// close() cancels their write scope. BLE's retry path backs off 500ms, so this window
T8b949e// also covers one retry on flaky GATT links.
Tff7b72try Tb4b4b4{
Tff7b72if Tb4b4b4(
Te6edf3sendPoliteDisconnect Tff7b72&Tff7b72&
Te6edf3currentTransport Tff7b72!Tff7b72= Tff7b72null Tff7b72&Tff7b72&
Te6edf3_connectionStateTb4b4b4.Te6edf3value Tff7b72!Tff7b72= Te6edf3ConnectionStateTb4b4b4.Te6edf3Disconnected
Tb4b4b4) Tb4b4b4{
Te6edf3isStopping Tff7b72= Tff7b72true
Te6edf3ignoreExceptionSuspend Tb4b4b4{
Te6edf3currentTransportTb4b4b4.Te6edf3handleSendToRadioTb4b4b4(Te6edf3ToRadioTb4b4b4(Te6edf3disconnect Tff7b72= Tff7b72trueTb4b4b4)Tb4b4b4.Te6edf3encodeTb4b4b4(Tb4b4b4)Tb4b4b4)
Te6edf3delayTb4b4b4(Te6edf3POLITE_DISCONNECT_DRAIN_MSTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4} Tff7b72finally Tb4b4b4{
Te6edf3isStarted Tff7b72= Tff7b72false
Te6edf3radioTransport Tff7b72= Tff7b72null
Te6edf3runningTransportId Tff7b72= Tff7b72null
Te6edf3isStopping Tff7b72= Tff7b72false
Tff7b72try Tb4b4b4{
Te6edf3currentTransportTff7b72?.Te6edf3closeTb4b4b4(Tb4b4b4)
Tb4b4b4} Tff7b72finally Tb4b4b4{
Te6edf3_serviceScopeTb4b4b4.Te6edf3cancelTb4b4b4(Ta5d6ff"Ta5d6ffstopping transportTa5d6ff"Tb4b4b4)
Te6edf3_serviceScope Tff7b72= Te6edf3CoroutineScopeTb4b4b4(Te6edf3dispatchersTb4b4b4.Te6edf3io Tff7b72+ Te6edf3SupervisorJobTb4b4b4(Tb4b4b4)Tb4b4b4)
Tff7b72if Tb4b4b4(Te6edf3notifyPermanent Tff7b72&Tff7b72& Te6edf3currentTransport Tff7b72!Tff7b72= Tff7b72nullTb4b4b4) Tb4b4b4{
Te6edf3onDisconnectTb4b4b4(Te6edf3isPermanent Tff7b72= Tff7b72trueTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffstartHeartbeatTb4b4b4(Tb4b4b4) Tb4b4b4{
Te6edf3heartbeatJobTff7b72?.Te6edf3cancelTb4b4b4(Tb4b4b4)
Te6edf3lastDataReceivedMillis Tff7b72= Te6edf3nowTb4b4b4(Tb4b4b4)
Te6edf3heartbeatJob Tff7b72=
Te6edf3serviceScopeTb4b4b4.Te6edf3launch Tb4b4b4{
Tff7b72while Tb4b4b4(Tff7b72trueTb4b4b4) Tb4b4b4{
Te6edf3delayTb4b4b4(Te6edf3HEARTBEAT_INTERVAL_MILLISTb4b4b4)
Te6edf3keepAliveTb4b4b4(Tb4b4b4)
Te6edf3checkLivenessTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Detects zombie connections where the BLE stack didn't report a disconnect.
*
* If we believe we're connected but haven't received any data from the radio within [LIVENESS_TIMEOUT_MILLIS], the
* connection is likely dead. Signal a non-permanent disconnect so the reconnect machinery can take over.
*
* Uses [clockMillis] for the current time so tests can inject a deterministic clock.
*/
Tff7b72internal Tff7b72fun Td2a8ffcheckLivenessTb4b4b4(Tb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3_connectionStateTb4b4b4.Te6edf3value Tff7b72!Tff7b72= Te6edf3ConnectionStateTb4b4b4.Te6edf3ConnectedTb4b4b4) Tff7b72return
Tff7b72val Te6edf3silenceMs Tff7b72= Te6edf3nowTb4b4b4(Tb4b4b4) Tff7b72- Te6edf3lastDataReceivedMillis
Tff7b72if Tb4b4b4(Te6edf3silenceMs Tff7b72> Te6edf3LIVENESS_TIMEOUT_MILLISTb4b4b4) Tb4b4b4{
T8b949e// "Silence" = lastDataReceivedMillis not updated by handleFromRadio (no inbound
T8b949e// packets). Only BLE suffers from silent zombie sessions (no disconnect signal from
T8b949e// stack), so the liveness-restart path is BLE-only. For non-BLE transports we return
T8b949e// WITHOUT emitting a disconnect or mutating ConnectionState β there is no
T8b949e// transport-level timeout contract proving that silence past this threshold means
T8b949e// the session is dead.
Tff7b72if Tb4b4b4(Te6edf3runningTransportId Tff7b72!Tff7b72= Te6edf3InterfaceIdTb4b4b4.Te6edf3BLUETOOTHTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffIgnoring liveness timeout for non-BLE transport (silence: Tffd700${Te6edf3silenceMsTffd700}Ta5d6ffms)Ta5d6ff" Tb4b4b4}
Tff7b72return
Tb4b4b4}
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{
Ta5d6ff"Ta5d6ffLiveness check failed: no data received for Tffd700${Te6edf3silenceMsTffd700}Ta5d6ffms Ta5d6ff" Tff7b72+
Ta5d6ff"Ta5d6ff(threshold: Tffd700${Te6edf3LIVENESS_TIMEOUT_MILLISTffd700}Ta5d6ffms). Restarting BLE transport.Ta5d6ff"
Tb4b4b4}
T8b949e// Force transport restart to recover from silent zombie sessions where the BLE stack
T8b949e// did not report a disconnect. Uses the same processLifecycle scope and transportMutex
T8b949e// pattern as setDeviceAddress() to guarantee clean teardown/restart sequencing.
T8b949e// The onDisconnect notification is emitted INSIDE the compareAndSet guard so that a
T8b949e// double liveness-timeout (timer not cancelled between fires) does not produce
T8b949e// duplicate disconnect notifications for a single restart cycle.
Tff7b72if Tb4b4b4(Te6edf3isRestartingTb4b4b4.Te6edf3compareAndSetTb4b4b4(Te6edf3expect Tff7b72= Tff7b72falseTb4b4b4, Te6edf3update Tff7b72= Tff7b72trueTb4b4b4)Tb4b4b4) Tb4b4b4{
T8b949e// Silent recovery: emit the non-permanent state transition (DeviceSleep) so the
T8b949e// reconnect machinery takes over, but do NOT pass an errorMessage. Automatic
T8b949e// liveness recovery is self-healing β surfacing a modal dialog for a transient
T8b949e// condition the app already handled is confusing UX. The warning log above
T8b949e// remains the observability surface for this event.
T8b949e//
T8b949e// Ordering note (pre-existing): onDisconnect fires here, BEFORE the launched
T8b949e// restart coroutine's `connectionRequested` check below. If an explicit disconnect()
T8b949e// races this timeout, a spurious DeviceSleep emission can leak to observers. The
T8b949e// connectionRequested gate still prevents the worse outcome β transport resurrection
T8b949e// β so this is a benign UI-level transient, not a state-machine bug.
Te6edf3onDisconnectTb4b4b4(Te6edf3isPermanent Tff7b72= Tff7b72falseTb4b4b4)
Te6edf3processLifecycleTb4b4b4.Te6edf3coroutineScopeTb4b4b4.Te6edf3launch Tb4b4b4{
Tff7b72try Tb4b4b4{
Te6edf3transportMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
T8b949e// Defense against a race between checkLiveness() firing and a
T8b949e// concurrent disconnect(): if the user has torn the connection down
T8b949e// since the heartbeat scheduled this restart, leave it down. The
T8b949e// transport is already null after disconnect()'s stopTransportLocked().
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3connectionRequestedTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffSkipping liveness restart: connection no longer requestedTa5d6ff" Tb4b4b4}
Tff7b72returnTf0883e@withLock
Tb4b4b4}
Te6edf3ignoreExceptionSuspend Tb4b4b4{
Te6edf3stopTransportLockedTb4b4b4(Te6edf3notifyPermanent Tff7b72= Tff7b72falseTb4b4b4, Te6edf3sendPoliteDisconnect Tff7b72= Tff7b72falseTb4b4b4)
Tb4b4b4}
T8b949e// This restart is launched from a heartbeat callback and has no caller to receive a
T8b949e// recoverable factory exception. The start path rolls back the failed session first;
T8b949e// contain the rethrow here so a failed reconnect does not crash the app.
Te6edf3ignoreExceptionSuspend Tb4b4b4{ Te6edf3startTransportLockedTb4b4b4(Tb4b4b4) Tb4b4b4}
Tb4b4b4}
Tb4b4b4} Tff7b72finally Tb4b4b4{
Te6edf3isRestartingTb4b4b4.Te6edf3value Tff7b72= Tff7b72false
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72fun Td2a8ffkeepAliveTb4b4b4(Te6edf3nowTb4b4b4: Tffa657Long Tff7b72= Te6edf3nowTb4b4b4(Tb4b4b4)Tb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3now Tff7b72- Te6edf3lastHeartbeatMillis Tff7b72> Te6edf3HEARTBEAT_INTERVAL_MILLISTb4b4b4) Tb4b4b4{
Te6edf3radioTransportTff7b72?.Te6edf3keepAliveTb4b4b4(Tb4b4b4)
Te6edf3lastHeartbeatMillis Tff7b72= Te6edf3now
Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffsendToRadioTb4b4b4(Te6edf3bytesTb4b4b4: Te6edf3ByteArrayTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3isStoppingTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffsendToRadio: transport stopping, dropping Tffd700${Te6edf3bytesTb4b4b4.Te6edf3sizeTffd700}Ta5d6ff bytesTa5d6ff" Tb4b4b4}
Tff7b72return
Tb4b4b4}
T8b949e// Snapshot the transport to avoid calling handleSendToRadio on a null reference.
T8b949e// There is still a benign race: stopTransportLocked() may cancel _serviceScope
T8b949e// between the null-check and the launch, causing the coroutine to be silently
T8b949e// dropped. This is acceptable β if the transport is shutting down, dropping the
T8b949e// send is the correct behavior.
Tff7b72val Te6edf3currentTransport Tff7b72=
Te6edf3radioTransport
Tff7b72?: Te6edf3run Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffsendToRadio: no active radio transport, dropping Tffd700${Te6edf3bytesTb4b4b4.Te6edf3sizeTffd700}Ta5d6ff bytesTa5d6ff" Tb4b4b4}
Tff7b72return
Tb4b4b4}
Te6edf3_serviceScopeTb4b4b4.Te6edf3handledLaunch Tb4b4b4{
Te6edf3currentTransportTb4b4b4.Te6edf3handleSendToRadioTb4b4b4(Te6edf3bytesTb4b4b4)
Te6edf3_meshActivityTb4b4b4.Te6edf3tryEmitTb4b4b4(Te6edf3MeshActivityTb4b4b4.Te6edf3SendTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4)
Tff7b72override Tff7b72fun Td2a8ffhandleFromRadioTb4b4b4(Te6edf3bytesTb4b4b4: Te6edf3ByteArrayTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3admitted Tff7b72=
Te6edf3synchronizedTb4b4b4(Te6edf3sessionCallbackLockTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3session Tff7b72= Te6edf3activeTransportSession Tff7b72?: Tff7b72returnTf0883e@synchronized Tff7b72false
Te6edf3enqueueReceivedDataTb4b4b4(Te6edf3bytesTb4b4b4, Te6edf3sessionTb4b4b4)
Tff7b72true
Tb4b4b4}
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3admittedTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffDropping Tffd700${Te6edf3bytesTb4b4b4.Te6edf3sizeTffd700}Ta5d6ff received bytes without an active transport sessionTa5d6ff" Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4)
Tff7b72private Tff7b72fun Td2a8ffenqueueReceivedDataTb4b4b4(Te6edf3bytesTb4b4b4: Te6edf3ByteArrayTb4b4b4, Te6edf3sessionTb4b4b4: Te6edf3RadioTransportSessionTb4b4b4) Tb4b4b4{
Tff7b72try Tb4b4b4{
Te6edf3lastDataReceivedMillis Tff7b72= Te6edf3nowTb4b4b4(Tb4b4b4)
T8b949e// trySend synchronously onto the Channel so packet order matches arrival order. The
T8b949e// previous `launch { emit() }` pattern dispatched each packet onto a fresh coroutine,
T8b949e// letting the scheduler reorder them β which broke the firmware config handshake
T8b949e// (see PhoneAPI.cpp initial-handshake sequence).
T8b949e// Reject before the copy: the channel bounds the frame COUNT, so without a per-frame ceiling a transport
T8b949e// handing over an oversized buffer defeats the memory bound the capacity is supposed to give. The stream
T8b949e// codec already enforces this on its own path; BLE passes through whatever the GATT read returned.
Tff7b72if Tb4b4b4(Te6edf3bytesTb4b4b4.Te6edf3size Tff7b72> Te6edf3MAX_FRAME_BYTESTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffDiscarding oversized Tffd700${Te6edf3bytesTb4b4b4.Te6edf3sizeTffd700}Ta5d6ff-byte frame (max Tffd700$Te6edf3MAX_FRAME_BYTESTa5d6ff)Ta5d6ff" Tb4b4b4}
Tff7b72return
Tb4b4b4}
Tff7b72val Te6edf3frame Tff7b72= Te6edf3ReceivedRadioFrameTb4b4b4(Te6edf3payload Tff7b72= Te6edf3bytesTb4b4b4.Te6edf3toByteStringTb4b4b4(Tb4b4b4)Tb4b4b4, Te6edf3session Tff7b72= Te6edf3sessionTb4b4b4.Te6edf3contextTb4b4b4)
Tff7b72val Te6edf3result Tff7b72= Te6edf3_receivedDataTb4b4b4.Te6edf3trySendTb4b4b4(Te6edf3frameTb4b4b4)
Tff7b72if Tb4b4b4(Te6edf3resultTb4b4b4.Te6edf3isFailureTb4b4b4) Tb4b4b4{
T8b949e// Rate-limited on purpose: drops only happen under sustained inbound traffic, and Kermit forwards to
T8b949e// Datadog/Crashlytics, so logging every drop would turn a bounded memory problem into unbounded
T8b949e// network and battery use. The counter is deliberately unsynchronised β a racy count only skews a
T8b949e// diagnostic line. A full queue reports failure with no exception; a closed channel carries one.
Tff7b72val Te6edf3drops Tff7b72= Tff7b72+Tff7b72+Te6edf3droppedFrameCount
Tff7b72if Tb4b4b4(Te6edf3drops Tff7b72=Tff7b72= T79c0ff1L Tff7b72|Tff7b72| Te6edf3drops Tff7b72% Te6edf3DROP_LOG_INTERVAL Tff7b72=Tff7b72= T79c0ff0LTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3resultTb4b4b4.Te6edf3exceptionOrNullTb4b4b4(Tb4b4b4)Tb4b4b4) Tb4b4b4{
Ta5d6ff"Ta5d6ffDropped Tffd700${Te6edf3bytesTb4b4b4.Te6edf3sizeTffd700}Ta5d6ff received bytes (Tffd700$Te6edf3dropsTa5d6ff total); receive queue at capacity Ta5d6ff" Tff7b72+
Ta5d6ff"Tffd700$Te6edf3RECEIVE_QUEUE_CAPACITYTa5d6ff or closedTa5d6ff"
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Te6edf3_meshActivityTb4b4b4.Te6edf3tryEmitTb4b4b4(Te6edf3MeshActivityTb4b4b4.Te6edf3ReceiveTb4b4b4)
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3tTb4b4b4: Te6edf3ThrowableTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3eTb4b4b4(Te6edf3tTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffhandleFromRadio failed while emitting dataTa5d6ff" Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffresetReceivedBufferTb4b4b4(Tb4b4b4) Tb4b4b4{
T8b949e// Drain any bytes buffered while no collector was attached. Without this, a stop/start cycle
T8b949e// would replay stale bytes ahead of the next session's firmware handshake, since the channel
T8b949e// outlives the orchestrator's per-start scope.
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffEmptyWhileBlockTa5d6ff"Tb4b4b4, Ta5d6ff"Ta5d6ffControlFlowWithEmptyBodyTa5d6ff"Tb4b4b4)
Tff7b72while Tb4b4b4(Te6edf3_receivedDataTb4b4b4.Te6edf3tryReceiveTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3isSuccessTb4b4b4) Tb4b4b4{Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffonConnectTb4b4b4(Tb4b4b4) Tb4b4b4{
Te6edf3synchronizedTb4b4b4(Te6edf3sessionCallbackLockTb4b4b4) Tb4b4b4{ Te6edf3publishConnectedTb4b4b4(Tb4b4b4) Tb4b4b4}
Tb4b4b4}
T8b949e/** Applies a callback already admitted under [sessionCallbackLock]. */
Tff7b72private Tff7b72fun Td2a8ffpublishConnectedTb4b4b4(Tb4b4b4) Tb4b4b4{
T8b949e// MutableStateFlow.value is thread-safe (backed by atomics) β assign directly rather than
T8b949e// launching a coroutine. The async launch pattern introduced a window where a concurrent
T8b949e// onDisconnect launch could execute AFTER an onConnect launch, leaving the service stuck
T8b949e// in Connected while the transport was actually disconnected.
Te6edf3lastDataReceivedMillis Tff7b72= Te6edf3nowTb4b4b4(Tb4b4b4)
Tff7b72if Tb4b4b4(Te6edf3_connectionStateTb4b4b4.Te6edf3value Tff7b72!Tff7b72= Te6edf3ConnectionStateTb4b4b4.Te6edf3ConnectedTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffBroadcasting connection state change to ConnectedTa5d6ff" Tb4b4b4}
Te6edf3_connectionStateTb4b4b4.Te6edf3value Tff7b72= Te6edf3ConnectionStateTb4b4b4.Te6edf3Connected
Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffonDisconnectTb4b4b4(Te6edf3isPermanentTb4b4b4: Tffa657BooleanTb4b4b4, Te6edf3errorMessageTb4b4b4: Tffa657String?Tb4b4b4, Te6edf3reasonTb4b4b4: Te6edf3TransportDisconnectReason?Tb4b4b4) Tb4b4b4{
Te6edf3synchronizedTb4b4b4(Te6edf3sessionCallbackLockTb4b4b4) Tb4b4b4{ Te6edf3publishDisconnectedTb4b4b4(Te6edf3isPermanentTb4b4b4, Te6edf3errorMessageTb4b4b4, Te6edf3reasonTb4b4b4) Tb4b4b4}
Tb4b4b4}
T8b949e/** Applies a callback already admitted under [sessionCallbackLock]. */
Tff7b72private Tff7b72fun Td2a8ffpublishDisconnectedTb4b4b4(Te6edf3isPermanentTb4b4b4: Tffa657BooleanTb4b4b4, Te6edf3errorMessageTb4b4b4: Tffa657String?Tb4b4b4, Te6edf3reasonTb4b4b4: Te6edf3TransportDisconnectReason?Tb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3resolvedErrorMessage Tff7b72= Te6edf3errorMessage Tff7b72?: Te6edf3reasonTff7b72?.Te6edf3toConnectionErrorMessageTb4b4b4(Tb4b4b4)
Tff7b72if Tb4b4b4(Te6edf3resolvedErrorMessage Tff7b72!Tff7b72= Tff7b72nullTb4b4b4) Tb4b4b4{
Te6edf3processLifecycleTb4b4b4.Te6edf3coroutineScopeTb4b4b4.Te6edf3launchTb4b4b4(Te6edf3dispatchersTb4b4b4.Te6edf3defaultTb4b4b4) Tb4b4b4{ Te6edf3_connectionErrorTb4b4b4.Te6edf3emitTb4b4b4(Te6edf3resolvedErrorMessageTb4b4b4) Tb4b4b4}
Tb4b4b4}
Tff7b72val Te6edf3newTargetState Tff7b72= Tff7b72if Tb4b4b4(Te6edf3isPermanentTb4b4b4) Te6edf3ConnectionStateTb4b4b4.Te6edf3Disconnected Tff7b72else Te6edf3ConnectionStateTb4b4b4.Te6edf3DeviceSleep
Tff7b72if Tb4b4b4(Te6edf3_connectionStateTb4b4b4.Te6edf3value Tff7b72!Tff7b72= Te6edf3newTargetStateTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffBroadcasting connection state change to Tffd700$Te6edf3newTargetStateTa5d6ff" Tb4b4b4}
Te6edf3_connectionStateTb4b4b4.Te6edf3value Tff7b72= Te6edf3newTargetState
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Per-transport-session wrapper around this service. Delegates every [RadioInterfaceService] surface (scope,
* address, sendToRadio, etc.) to the enclosing service; only the three [RadioTransportCallback] entry points are
* gated on the captured [session] still being the [activeTransportSession]. Admission and the synchronous callback
* side effect share [sessionCallbackLock] with teardown, eliminating a check-then-use window. Late callbacks from a
* torn-down transport are dropped BEFORE bytes enter the shared received-data channel or connection/config state
* mutates. Address is logged with [anonymize] and the session generation only β never the raw address bytes.
*/
Tff7b72private Tff7b72inner Tff7b72class T56d364SessionBoundRadioInterfaceServiceTb4b4b4(Tff7b72val Te6edf3sessionTb4b4b4: Te6edf3RadioTransportSessionTb4b4b4) Tb4b4b4:
Te6edf3RadioInterfaceService Tff7b72by Tff7b72thisTf0883e@SharedRadioInterfaceService Tb4b4b4{
Tff7b72override Tff7b72fun Td2a8ffonConnectTb4b4b4(Tb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3admitted Tff7b72= Te6edf3runIfTransportSessionActiveTb4b4b4(Te6edf3sessionTb4b4b4, Tff7b72::Te6edf3publishConnectedTb4b4b4)
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3admittedTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffDropping stale onConnect gen=Tffd700${Te6edf3sessionTb4b4b4.Te6edf3generationTffd700}Ta5d6ff addr=Tffd700${Te6edf3sessionTb4b4b4.Te6edf3addressTb4b4b4.Te6edf3anonymizeTffd700}Ta5d6ff" Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffonDisconnectTb4b4b4(Te6edf3isPermanentTb4b4b4: Tffa657BooleanTb4b4b4, Te6edf3errorMessageTb4b4b4: Tffa657String?Tb4b4b4, Te6edf3reasonTb4b4b4: Te6edf3TransportDisconnectReason?Tb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3admitted Tff7b72=
Te6edf3runIfTransportSessionActiveTb4b4b4(Te6edf3sessionTb4b4b4) Tb4b4b4{ Te6edf3publishDisconnectedTb4b4b4(Te6edf3isPermanentTb4b4b4, Te6edf3errorMessageTb4b4b4, Te6edf3reasonTb4b4b4) Tb4b4b4}
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3admittedTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffDropping stale onDisconnect gen=Tffd700${Te6edf3sessionTb4b4b4.Te6edf3generationTffd700}Ta5d6ff addr=Tffd700${Te6edf3sessionTb4b4b4.Te6edf3addressTb4b4b4.Te6edf3anonymizeTffd700}Ta5d6ff" Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffhandleFromRadioTb4b4b4(Te6edf3bytesTb4b4b4: Te6edf3ByteArrayTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3admitted Tff7b72=
Te6edf3runIfTransportSessionActiveTb4b4b4(Te6edf3sessionTb4b4b4) Tb4b4b4{
Tff7b72thisTf0883e@SharedRadioInterfaceService.enqueueReceivedDataTb4b4b4(Te6edf3bytesTb4b4b4, Te6edf3sessionTb4b4b4)
Tb4b4b4}
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3admittedTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{
Ta5d6ff"Ta5d6ffDropping stale handleFromRadio (Tffd700${Te6edf3bytesTb4b4b4.Te6edf3sizeTffd700}Ta5d6ff bytes) gen=Tffd700${Te6edf3sessionTb4b4b4.Te6edf3generationTffd700}Ta5d6ff Ta5d6ff" Tff7b72+
Ta5d6ff"Ta5d6ffaddr=Tffd700${Te6edf3sessionTb4b4b4.Te6edf3addressTb4b4b4.Te6edf3anonymizeTffd700}Ta5d6ff"
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Served by rngit 1.5.0 - Generated in 0.19s